Python
Python 多线程与多进程:GIL、并发模型与实践
从 GIL 出发比较 threading、multiprocessing 与 concurrent.futures 的适用场景。

ABSTRACT 概念速览
- 多线程(threading):多个线程共享同一进程的内存空间;受 GIL 限制,适合 I/O 密集型 任务
- 多进程(multiprocessing):每个进程有独立内存空间;绕过 GIL,适合 CPU 密集型 任务
一、GIL —— 理解并发的前提
NOTE 什么是 GIL? Global Interpreter Lock(全局解释器锁),CPython 的实现机制:同一时刻只有一个线程能执行 Python 字节码。
- I/O 操作(网络、磁盘)会释放 GIL → 多线程有效
- 纯计算不释放 GIL → 多线程无法真正并行,用多进程代替
| 场景 | 推荐方案 | 原因 |
|---|---|---|
| 网络请求、文件读写 | threading / asyncio |
I/O 等待期间释放 GIL |
| 图像处理、数学计算 | multiprocessing |
独立进程,真正并行 |
| 两者都有 | concurrent.futures |
统一接口,灵活切换 |
二、threading —— 多线程
2.1 创建线程
import threading
def task(name, count):
for i in range(count):
print(f"[{name}] step {i}")
# 方式一:直接传函数
t = threading.Thread(target=task, args=("A", 3), kwargs={})
t.start() # 启动线程
t.join() # 等待线程结束(阻塞主线程)
# 方式二:继承 Thread 类(适合逻辑复杂时)
class MyThread(threading.Thread):
def __init__(self, name):
super().__init__()
self.name = name
def run(self): # 重写 run(),不要直接调用 run()
print(f"{self.name} running")
t = MyThread("Worker")
t.start()
t.join()TIP
start()vsrun()
t.start()→ 在新线程中执行run()t.run()→ 在当前线程直接调用,没有新线程!
2.2 常用属性与方法
t = threading.Thread(target=task, args=("A",), daemon=True)
t.start() # 启动
t.join(timeout=5.0) # 最多等待 5 秒
t.is_alive() # 是否还在运行 → bool
t.name # 线程名(可自定义)
t.daemon # 守护线程:主线程结束时强制退出
threading.current_thread() # 返回当前线程对象
threading.enumerate() # 返回所有存活线程列表
threading.active_count() # 存活线程数量WARNING 守护线程(daemon)
daemon=True时,主线程结束,守护线程立即被杀死,不等其完成。
适合后台心跳、日志刷新等"可抛弃"任务。
2.3 线程锁(Lock)
多线程共享内存 → 需要加锁防止竞态条件:
import threading
counter = 0
lock = threading.Lock()
def increment():
global counter
with lock: # 推荐:自动 acquire / release
counter += 1
# 等价于:
def increment_manual():
global counter
lock.acquire()
try:
counter += 1
finally:
lock.release() # 必须放在 finally!
threads = [threading.Thread(target=increment) for _ in range(1000)]
for t in threads: t.start()
for t in threads: t.join()
print(counter) # 1000(无锁时结果不确定)2.4 RLock —— 可重入锁
rlock = threading.RLock()
def outer():
with rlock:
inner() # 同一线程可再次获取,不会死锁
def inner():
with rlock: # Lock 在这里会死锁,RLock 不会
print("inner")2.5 Semaphore —— 信号量(限制并发数)
sem = threading.Semaphore(3) # 最多 3 个线程同时运行
def limited_task(i):
with sem:
print(f"Task {i} running")
time.sleep(1)
threads = [threading.Thread(target=limited_task, args=(i,)) for i in range(10)]
for t in threads: t.start()
for t in threads: t.join()2.6 Event —— 线程间通信
event = threading.Event()
def waiter():
print("等待信号...")
event.wait() # 阻塞,直到 event 被 set()
print("收到信号,继续!")
def setter():
time.sleep(2)
event.set() # 通知所有等待的线程
threading.Thread(target=waiter).start()
threading.Thread(target=setter).start()| 方法 | 说明 |
|---|---|
event.set() |
设置标志,释放所有 wait() |
event.clear() |
清除标志,后续 wait() 会再次阻塞 |
event.is_set() |
返回当前状态 |
event.wait(timeout) |
阻塞等待,可设超时 |
2.7 Queue —— 线程安全队列
from queue import Queue
q = Queue(maxsize=5) # 0 表示无限制
q.put(item) # 放入,队满则阻塞
q.put_nowait(item) # 队满直接抛 Full 异常
q.get() # 取出,队空则阻塞
q.get_nowait() # 队空直接抛 Empty 异常
q.task_done() # 告知该任务已完成
q.join() # 阻塞,直到所有 task_done() 被调用
q.qsize() # 当前大小
q.empty() / q.full() # 状态判断(非原子,仅供参考)生产者-消费者模式:
from queue import Queue
import threading, time
q = Queue()
def producer():
for i in range(5):
q.put(i)
print(f"生产: {i}")
time.sleep(0.5)
def consumer():
while True:
item = q.get()
if item is None: # 哨兵值:退出信号
break
print(f"消费: {item}")
q.task_done()
t1 = threading.Thread(target=producer)
t2 = threading.Thread(target=consumer)
t1.start(); t2.start()
t1.join()
q.put(None) # 发送退出信号
t2.join()三、multiprocessing —— 多进程
3.1 创建进程
from multiprocessing import Process
def task(name):
print(f"进程 {name} 运行中,PID={os.getpid()}")
if __name__ == "__main__": # ⚠️ Windows 必须有这个守卫!
p = Process(target=task, args=("Worker",))
p.start()
p.join()DANGER
if __name__ == "__main__"守卫 Windows 使用spawn方式创建进程,会重新 import 主模块。
没有这个守卫会无限递归创建进程! Linux/macOS 用fork,风险较小但仍建议加上。
3.2 常用属性与方法
p = Process(target=task, args=("A",), daemon=False)
p.start()
p.join(timeout=5)
p.terminate() # 发送 SIGTERM(较优雅)
p.kill() # 发送 SIGKILL(立即杀死)
p.is_alive()
p.pid # 进程 ID
p.exitcode # 退出码(None 表示还在运行)3.3 进程间通信(IPC)
进程内存相互隔离,需要专门的 IPC 机制:
Queue(多进程安全)
from multiprocessing import Process, Queue
def producer(q):
q.put("hello from producer")
def consumer(q):
msg = q.get()
print(f"收到: {msg}")
if __name__ == "__main__":
q = Queue()
p1 = Process(target=producer, args=(q,))
p2 = Process(target=consumer, args=(q,))
p1.start(); p2.start()
p1.join(); p2.join()Pipe(双向管道)
from multiprocessing import Process, Pipe
def child(conn):
conn.send("来自子进程")
msg = conn.recv()
print(f"子进程收到: {msg}")
conn.close()
if __name__ == "__main__":
parent_conn, child_conn = Pipe() # 返回两端
p = Process(target=child, args=(child_conn,))
p.start()
print(f"父进程收到: {parent_conn.recv()}")
parent_conn.send("来自父进程")
p.join()共享内存(Value / Array)
from multiprocessing import Process, Value, Array
import ctypes
def modify(val, arr):
val.value = 42
arr[0] = 100
if __name__ == "__main__":
v = Value(ctypes.c_int, 0) # 共享整数
a = Array(ctypes.c_int, [1,2,3]) # 共享数组
p = Process(target=modify, args=(v, a))
p.start()
p.join()
print(v.value, list(a)) # 42 [100, 2, 3]NOTE 共享内存 vs Queue/Pipe
Value/Array:零拷贝,高性能,适合大量数值数据Queue/Pipe:序列化传输,更灵活,适合任意 Python 对象
3.4 进程锁
from multiprocessing import Process, Lock
def task(lock, i):
with lock:
print(f"进程 {i} 安全写入")
if __name__ == "__main__":
lock = Lock()
processes = [Process(target=task, args=(lock, i)) for i in range(5)]
for p in processes: p.start()
for p in processes: p.join()3.5 进程池(Pool)
TIP 何时用 Pool? 需要执行大量同类任务时,Pool 自动管理进程复用,避免频繁创建销毁进程的开销。
from multiprocessing import Pool
def square(x):
return x * x
if __name__ == "__main__":
with Pool(processes=4) as pool: # 4 个工作进程
# map:阻塞,返回有序结果列表
result = pool.map(square, range(10))
print(result) # [0, 1, 4, 9, 16, 25, 36, 49, 64, 81]
# apply:同步执行单个任务
r = pool.apply(square, args=(5,))
print(r) # 25
# apply_async:异步执行单个任务
async_result = pool.apply_async(square, args=(7,))
print(async_result.get(timeout=5)) # 49
# map_async:异步 map
async_results = pool.map_async(square, range(5))
print(async_results.get()) # [0, 1, 4, 9, 16]
# imap:惰性迭代(大数据节省内存)
for r in pool.imap(square, range(5)):
print(r)| 方法 | 是否阻塞 | 返回值 | 适用场景 |
|---|---|---|---|
map(f, iterable) |
✅ 阻塞 | list |
同步批量处理 |
map_async(f, iterable) |
❌ 非阻塞 | AsyncResult |
异步批量处理 |
apply(f, args) |
✅ 阻塞 | 单个结果 | 同步单次调用 |
apply_async(f, args) |
❌ 非阻塞 | AsyncResult |
异步单次调用 |
imap(f, iterable) |
✅ 惰性 | 迭代器 | 大数据流式处理 |
starmap(f, iterable) |
✅ 阻塞 | list |
多参数展开 |
# starmap:多参数
def add(a, b): return a + b
with Pool(4) as pool:
result = pool.starmap(add, [(1,2), (3,4), (5,6)])
print(result) # [3, 7, 11]四、concurrent.futures —— 统一高层接口
ABSTRACT 核心优势 用同一套 API 切换多线程和多进程,代码改动极小。
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
def task(n):
return n * n
# 多线程(I/O 密集型)
with ThreadPoolExecutor(max_workers=4) as executor:
futures = [executor.submit(task, i) for i in range(5)]
for f in futures:
print(f.result()) # 阻塞等待结果
# 多进程(CPU 密集型)—— 只改一行!
with ProcessPoolExecutor(max_workers=4) as executor:
futures = [executor.submit(task, i) for i in range(5)]
for f in futures:
print(f.result())map 方式
with ThreadPoolExecutor(max_workers=4) as executor:
# 按提交顺序返回结果
results = list(executor.map(task, range(10)))
print(results)Future 对象
future = executor.submit(task, 42)
future.result(timeout=5) # 获取结果,超时抛 TimeoutError
future.done() # 是否完成
future.cancelled() # 是否被取消
future.cancel() # 尝试取消(若已开始执行则无效)
future.exception() # 若任务异常,返回异常对象as_completed —— 哪个先完成先处理
from concurrent.futures import as_completed
import time, random
def slow_task(n):
time.sleep(random.random())
return n * 2
with ThreadPoolExecutor(max_workers=4) as executor:
futures = {executor.submit(slow_task, i): i for i in range(5)}
for future in as_completed(futures): # 先完成先返回
n = futures[future]
print(f"任务 {n} 完成,结果: {future.result()}")五、对比速查表
多线程 vs 多进程
| 维度 | threading |
multiprocessing |
|---|---|---|
| 内存 | 共享 | 独立 |
| GIL 影响 | 受限(I/O 可并发) | 不受限(真正并行) |
| 适合场景 | I/O 密集 | CPU 密集 |
| 通信方式 | 共享变量 + 锁 | Queue / Pipe / 共享内存 |
| 创建开销 | 小 | 大 |
| 稳定性 | 崩溃影响整个进程 | 进程间隔离 |
同步原语速查
| 原语 | 模块 | 作用 |
|---|---|---|
Lock |
threading / multiprocessing | 互斥锁,同一时刻只有一个持有者 |
RLock |
threading | 可重入锁,同一线程可多次获取 |
Semaphore |
threading / multiprocessing | 限制并发数量 |
Event |
threading / multiprocessing | 线程/进程间信号通知 |
Condition |
threading | 复杂条件同步(wait/notify) |
Barrier |
threading | 等待所有线程到达同一屏障点 |
六、常见陷阱
DANGER 陷阱 1:多进程忘写
__main__守卫# ❌ Windows 下无限递归 p = Process(target=f) p.start() # ✅ 正确 if __name__ == "__main__": p = Process(target=f) p.start()
DANGER 陷阱 2:线程中修改共享变量不加锁
# ❌ 竞态条件,结果不确定 counter += 1 # ✅ 加锁保护 with lock: counter += 1
DANGER 陷阱 3:多进程中普通全局变量不共享
count = 0 # 每个进程有自己的副本,修改互不影响 # ✅ 需要共享时用 Value / Manager from multiprocessing import Value count = Value('i', 0)
DANGER 陷阱 4:子线程异常不会传播到主线程
# ❌ 子线程抛异常,主线程不知道 t = Thread(target=risky_func) t.start(); t.join() # 不会抛异常! # ✅ 用 concurrent.futures 会把异常保存到 Future future = executor.submit(risky_func) future.result() # 这里才会重新抛出异常
WARNING 陷阱 5:死锁 多个锁按不同顺序获取时可能死锁:
# 线程 1: 先拿 lock_a 再拿 lock_b # 线程 2: 先拿 lock_b 再拿 lock_a # → 互相等待,永远阻塞 # ✅ 统一加锁顺序,或使用 timeout lock_a.acquire(timeout=1)
七、选型决策树
需要并发?
├─ I/O 密集(网络/文件)
│ ├─ 任务数不多 → threading.Thread
│ ├─ 任务数很多 → ThreadPoolExecutor
│ └─ 大量异步 I/O → asyncio(参见异步笔记)
│
└─ CPU 密集(计算)
├─ 任务数不多 → multiprocessing.Process
├─ 任务数很多 → ProcessPoolExecutor 或 Pool
└─ 需要共享大量数据 → multiprocessing + shared memory关联笔记:Python异步编程 · Python GIL深入 · asyncio事件循环